Conversation
…r-for-jvmdispatch
…r-for-jvmdispatch
…r-as-spark=memory-consumer-for-jvmdispatch
…r-for-jvmdispatch
…r-for-jvmdispatch
|
@sunchao would you like to take a look at this? TIA! |
sunchao
left a comment
There was a problem hiding this comment.
Two P2 findings in task-allocator concurrency, detailed inline. Validated with targeted JVM reproductions against this exact head and Arrow 18.3.0. The native producer/cleanup paths were source-traced, not run end-to-end.
|
@peterxcli can you check the CI failures? |
@sunchao I opened a fix for it |
|
Ah thanks! Approved the PR and will merge soon |
…r-for-jvmdispatch
…r-for-jvmdispatch
Native operators that retain an exported UDF result (hash join builds, shuffle buffers) take their own Spark reservations for those buffers, so keeping the JVM task charge until the FFI release charged the same physical buffers twice and could fail a native reservation that fit the pool. At export, transfer the result's buffer accounting from the task allocator to the root allocator (zero-copy; FFI addresses unchanged) and release the corresponding Spark charge. The UDF is still charged while it holds the memory; once native execution owns the buffers, the only Spark charge is whatever a retaining native operator reserves. Neither an ownership transfer nor the eventual FFI release fires the task allocator's AllocationListener, so the explicit release cannot double-free, and getAccountedSize() being non-zero only on the owning ledger keeps pass-through inputs and re-exported chunks at zero. The C schema is exported from the result's own Field: a transferred complex vector rebuilds child names from runtime data vectors, which Arrow hardcodes to "$data$". The C array carries no names, so it comes from the transferred vector. Exported buffers no longer pin TaskState past task completion, so the deferred-release retention (and the task-attempt-ID reuse concern it covered) is gone. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Getting JVM UDF output vectors onto Spark's task accounting is worth doing, and the distinction between imported native buffers (already accounted, keep on the root allocator) and JVM-created output buffers (charge to the task) is the right one. The retention logic that keeps the allocator alive until the last FFI release, while dropping Spark accounting at task completion, is carefully thought through. Four things. The threading contract for The class doc goes from:
to:
That is a breaking change to the contract that user-written Was the old guarantee wrong all along, or does this PR change when concurrent calls can happen? Either way this deserves its own line in the description and probably a migration-guide note, because it is a much bigger deal for users than the memory accounting is. Right now it reads like an incidental doc edit. Lock ordering deserves to be written down
I worked through it and did not find an actual inversion, but it took a while and I am not certain. Could you add a short comment naming the three locks and the order they must always be taken in? Anything that acquires a Spark memory lock from inside an Arrow allocation listener callback deserves that much. On-heap Tungsten mode silently does nothing this.consumer = taskMemoryManager.getTungstenMemoryMode() == MemoryMode.OFF_HEAP
? new TaskMemoryConsumer(taskMemoryManager) : null;With on-heap Tungsten memory the whole accounting path is disabled. That is probably correct, since Arrow's buffers are off-heap and charging them to Spark's on-heap pool would be wrong. But it is not stated anywhere, and a user running on-heap who reads the release notes will believe their UDF memory is accounted. Could the class doc say so explicitly, and ideally log once at debug level when accounting is skipped? The registration-order dependency
Is there a way to make that structural rather than conventional, for example having |
- Document the lock order (TaskMemoryManager monitor -> TaskState monitor -> MemoryManager monitor) on TaskState, verified against Spark 3.5.8 and 4.1.3: acquireExecutionMemory holds the TMM monitor across spills, releaseExecutionMemory never takes it, so every release-side path (TaskState -> MemoryManager) is a suffix of the acquire-side order and no inversion exists. - State explicitly that on-heap Tungsten memory disables Spark accounting for UDF Arrow buffers (they are off-heap, so there is no matching pool), and log that at debug level once per task. - Ground the CometUDF concurrency contract in its actual causes: DataFusion pipelining via spawned Tokio tasks (HashJoinExec OnceAsync) and multiple prefetching native plans per Spark task. CometScalaUDFCodegen has synchronized evaluate for this reason since its introduction; the old "at most one thread" trait doc was already contradicted on main. - Make completion-listener registration structural rather than order-dependent: register outside computeIfAbsent so a straggler evaluation racing task teardown recreates state safely instead of re-entering the map from its own mapping function, and clarify that listener ordering is for orderly shutdown, not memory safety. New test covers stragglers both during and after teardown. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
|
@andygrove Thanks — evaluated all four; addressed in 1dc8fc7. (1) Threading contract: this is a doc correction, not a behavior change. The "at most one thread" wording was already contradicted on main: (2) Lock ordering: added a javadoc on (3) On-heap mode: agreed — the behavior is correct (Arrow buffers are off-heap; there is no matching Spark pool), but it is now stated in the (4) Registration order: the failure mode is not use-after-free — the allocator only closes once in-flight evaluations finish and it holds no memory, and post-completion evaluation is rejected. But investigating this surfaced a real edge: a straggler |
sunchao
left a comment
There was a problem hiding this comment.
Two P2 findings in output-buffer accounting, detailed inline. The empty-output issue was reproduced through a native SQL query; the shared-scratch issue was reproduced through the public bridge with a custom CometUDF.
…F export charge chargedOutputSize computed the Spark charge to release at export from getBuffers(false), which omits the allocated buffers of zero-length children (an empty list's data vector keeps its allocated capacity) even though TransferPair.transfer() moves those chunks to the root allocator; the uncounted bytes stayed reserved against the task until completion. It now enumerates each vector's physical field buffers recursively, preserving per-ledger deduplication. It also released the full chunk charge for buffers shared with a retained scratch vector (an aligned splitAndTransfer slice shares the scratch ledger). When native execution released the FFI result, Arrow handed chunk ownership back to the surviving scratch ledger without a listener callback, and the eventual scratch close fired onRelease and freed the same charge a second time. A chunk is now counted only when every live reference to its ledger comes from the result tree; shared chunks keep their Spark charge and release it exactly once through onRelease, or wholesale at task completion if scratch closes while native still holds the buffers. Both regression tests fail against the previous logic: the empty-child test strands 32768 bytes of task charge at export, and the shared scratch test observes the charge dropped while a task ledger still references the chunk. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…r-for-jvmdispatch
…r-for-jvmdispatch
|
Triage note: #5998 overlaps with this. It attaches an They are not the same change. This PR cuts a per-task child allocator for the UDF output path specifically and charges it to a non-spillable |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: JVM UDF output allocations bypassed Spark task memory accounting.
- Design approach: Charge a task-scoped Arrow allocator through a non-spillable Spark consumer, then transfer output ownership at export. The transfer preserves shared buffers without copying their contents.
- Correctness / compatibility analysis: Found one introduced P1 compilation failure detailed below. Reviewed allocation, release and completion behavior against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources and Arrow 18.3.0. Earlier review fixes are present, and no remaining existing P1/P2 concern was substantiated.
- Key design decisions: Task-context identity separates task lifetimes. In-flight evaluation guards protect teardown. Recursive buffer accounting handles empty children and shared ledgers, keeping ownership logic within the bridge.
- Implementation sketch:
CometExecIteratorregisters task state, codegen receives the allocator throughCometUDF.evaluate, and the bridge transfers exported results to the root allocator. - Behavioral changes worth calling out: Off-heap UDF allocations now face Spark admission control. On-heap mode skips Spark charging. Custom implementations must adopt the new allocator parameter. Runtime performance was not measured.
- Suggested improvements: Update both benchmark callers for the new signature and rerun the affected build and tests.
Reviewed the entire 10-file diff from 36146a87bf9ca9ca9e211b4372628ed2f9d8c8c6 to dc117b0135659137591e26e84edd87ee1a387e03. Confirmed the PR is non-draft. Read existing reviews, issue comments, inline comments and threads. Routed skill: review-comet-pr. Checked sibling skill scopes; none applies to this allocator and lifecycle change.
Exact-head CI: All four Linux Spark 4.1 test groups and both TPC verification jobs fail during test compilation, making Required Checks fail. JVM/native builds, Rust tests and lint checks pass. The logs identify the merge of the requested head and base.
Validation limits: The local focused test attempt stopped before compilation because Maven Central access was unavailable. No local runtime tests or end-to-end queries ran. The finding is independently confirmed by exact-head CI logs and the unchanged benchmark callers. Project code remains unchanged.
|
|
||
| override def evaluate(inputs: Array[ValueVector], numRows: Int): ValueVector = { | ||
| override def evaluate( | ||
| allocator: BufferAllocator, |
There was a problem hiding this comment.
[P1] Update both benchmark callers when adding the allocator parameter. CometTimeExtractBenchmark.scala still calls dispatcher.evaluate(inputs, size) at lines 122 and 149. With the default Spark 4.1 profile, test compilation now fails instead of compiling those calls, preventing every Linux JVM test group and both TPC verification jobs from running. Please pass CometArrowAllocator at both sites, as the updated direct call in CometCodegenSuite does, then rerun test compilation and the focused suites.
Evidence: Exact-head CI run 35170049560 checks out the merge of dc117b0 into 36146a8. The expressions job reports not enough arguments for method evaluate and Unspecified value parameter numRows at CometTimeExtractBenchmark.scala:122 and :149, then fails scala-maven-plugin:4.9.6:testCompile. The same errors occur in the exec, scans, shuffle, TPC-H and TPC-DS jobs. Source comparison confirms the callers are unchanged while the two-argument method was replaced by this three-argument signature. https://github.com/apache/datafusion-comet/actions/runs/35170049560/job/105040966472
# Conflicts: # spark/src/main/java/org/apache/comet/udf/CometUdfBridge.java
|
I closed #5998 after finding that charging JVM Arrow memory to Spark, and refusing on a short grant, turns native spills into task failures. This PR refuses the same way:
UDF output goes straight into native operators that could have spilled instead, so I'd expect the same failure here. Could the consumer record the allocation without refusing it? Or refuse only when nothing native in the task can spill? For the rest of the JVM Arrow memory, I ended up with #6250: it reports the figures in the memory usage log and doesn't gate on them. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: JVM UDF output allocations bypassed Spark task memory accounting.
- Design approach: Introduce a task-scoped Arrow allocator backed by a non-spillable Spark consumer, then transfer output ownership at export.
- Correctness / compatibility analysis: No additional introduced P1/P2 issues found within this review. The earlier benchmark compilation blocker is fixed. The existing memory-pressure concern remains substantiated and unresolved, as detailed below.
- Key design decisions: Task-context identity separates task lifetimes. In-flight evaluation guards protect teardown. Recursive ledger accounting handles empty children and shared buffers. Ownership logic stays localized in the bridge.
- Implementation sketch:
CometExecIteratorregisters task state,CometUDF.evaluatereceives the allocator, and codegen allocates outputs through it. Export transfers buffers to the root allocator and releases exclusive output charges. - Behavioral changes worth calling out: Off-heap UDF allocations now face Spark admission control. On-heap mode skips Spark charging. Custom UDF implementations must adopt the allocator parameter. Fresh codegen outputs retain zero-copy export. Runtime overhead was not benchmarked.
- Suggested improvements: Resolve the existing interaction between UDF admission control and native spilling before merge, with a small-pool UDF-to-sort regression test.
Reviewed the entire 12-file diff from 81f2574d5e92028d6a0c8aebd39e6e0b5b4c2fd2 to dd24a38d428fd0cddc1a0cf1ce1601a4be91a3a3. Confirmed the PR is non-draft. Read AGENTS.md, existing reviews, issue comments, inline comments and review threads. Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-ffi-pr. Compared relevant behavior against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources, plus Arrow 18.3.0.
The existing memory-pressure concern remains a P2 blocker. An exact-source component probe reserved 67,104,768 bytes through CometTaskMemoryManager in a 64 MiB pool. A 1,024-row UDF requiring 8,192 Arrow bytes succeeded with the root-allocator control but failed with the task allocator. Partial grants were correctly returned, and the same allocation succeeded after releasing native reservations. Source tracing confirms NativeMemoryConsumer.spill() returns zero, while DataFusion’s sorter propagates an input error before reaching insert_batch and its spill-and-retry path. This can turn spillable workloads into task failures. It is already discussed, so no duplicate finding is added.
Exact-head CI: Run 36243579843 is green, including Required Checks, all four Linux Spark 4.1 test groups, native/Rust checks and TPC-H/TPC-DS verification. Logs confirm the merge includes the requested head. Spark SQL, macOS and other Spark runtime profiles were skipped.
Validation limits: All six CometUdfBridgeTest tests passed locally in an isolated build of the exact bridge sources with the original allocator declarations, Spark 4.1.3, Arrow 18.3.0 and JDK 21. The pressure reproduction was a JVM component probe, not an end-to-end native sort query. No full local project build or Spark SQL suite ran. Project files remain unchanged.
… grants The per-task UDF allocator's listener threw OutOfMemoryException from onPreAllocation whenever Spark's acquireMemory granted less than requested. Native operators reserve through a consumer whose spill returns 0 and spill only when their own try_grow fails, so under pressure native execution has already filled the pool and the next UDF output allocation is the one refused, failing the task where a downstream native sort or shuffle could have spilled. The listener now charges whatever Spark grants and carries the rest as a shortfall, like the native pools' grow/overcommit. Releases (buffer close, failed allocation, export to native) repay the shortfall before returning anything to Spark, so Spark stays charged for min(its grant, outstanding bytes) and is never handed back more than it granted. The consumer also clamps every release to what it still holds, which keeps Spark's books right when Arrow moves chunks between allocators without a listener callback. Task completion returns exactly the outstanding grant. The only remaining refusal is an allocation after task completion. Tests: the blocking-release test now requires the pending allocation to succeed; two JUnit tests reproduce the small-pool probe (native consumer holding all but 4 KiB of a 64 MiB pool) through the public bridge and on the release and completion paths; CometCodegenMemoryPressureSuite runs regexp_replace (codegen dispatch) into a native sort in a 16 MiB pool and asserts the sort spills. All fail on the previous head with "Failed to acquire ... JVM UDF bytes from Spark TaskMemoryManager".
|
This is a light fully automated review since there are so many PRs open.
|
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: JVM UDF Arrow allocations bypassed Spark task memory accounting.
- Design approach: Introduce a task-scoped allocator backed by a non-spillable Spark consumer. Record short grants without rejecting allocations, allowing downstream native operators to spill.
- Correctness / compatibility analysis: No additional introduced P1/P2 issues found within this review. Earlier compilation and memory-pressure blockers are fixed. The existing scratch-buffer retention concern remains a reproducible P2 blocker.
- Key design decisions: Task-context identity separates task lifetimes, in-flight guards protect teardown, and recursive ledger inspection handles empty children and shared buffers. Ownership logic remains centralized in the bridge.
- Implementation sketch:
CometExecIteratorregisters task state,CometUDF.evaluatereceives the allocator, and export transfers output ownership to the root allocator without copying buffer contents. - Behavioral changes worth calling out: Custom UDF implementations must adopt the allocator parameter. On-heap mode skips Spark charging. Short grants preserve execution but do not bound physical memory. Runtime accounting overhead was not benchmarked.
- Suggested improvements: Resolve the existing scratch-lifecycle concern by providing cleanup after in-flight evaluations finish, with a repeated-task regression test.
The existing scratch-buffer concern remains unresolved. A custom UDF retaining an 8,192-byte scratch buffer from the supplied allocator left one additional task state after each completed task. Three sequential tasks retained three child allocators and three completed task contexts. taskCompleted() clears UDF instances at CometUdfBridge.java:652, while closeIfIdle() requires an empty allocator at line 677. Consequently, the pre-existing scratch-buffer leak now also retains task state and its memory manager indefinitely. This is already reported, so no duplicate finding is added.
Reviewed the entire 16-file diff from 81f2574d5e92028d6a0c8aebd39e6e0b5b4c2fd2 to e8ed305d28643fe0fd13071d57e9ec4bf24351a5. Confirmed the PR remains non-draft. Read AGENTS.md, existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-ffi-pr, and review-comet-memory-pr. Compared relevant lifecycle and accounting semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources, plus Arrow 18.3.0.
Exact-head CI: Run 36468098758 passed, including all four Linux Spark 4.1 test groups, native/Rust checks, TPC verification and the new memory-pressure regression. Logs confirm the tested merge includes the requested head. A duplicate cancelled run produced a failed Required Checks result. Spark SQL, macOS and other Spark runtime profiles were skipped.
Validation limits: All eight CometUdfBridgeTest tests passed in a disposable component build using exact checkout sources and the original allocator declarations, with Spark 4.1.3, Arrow 18.3.0 and JDK 21. The scratch reproduction exercised the public JVM bridge and Arrow FFI, not a native query. No full local project build or Spark SQL suite ran. Project code and GitHub state remain unchanged.
Which issue does this PR close?
Closes #4174.
Rationale for this change
CometScalaUDFCodegen, currently the onlyCometUDFimplementation, creates output Arrow vectors from the process-wideCometArrowAllocator. These JVM allocations bypass Spark’s task memory accounting.FFI input vectors only wrap native-owned buffers that are already accounted by Comet, so they must not be charged again. Output buffers, however, are created on the JVM and may remain alive through the FFI release callback after the Java vector closes.
What changes are included in this PR?
MemoryConsumer.CometUDFinterface and codegen output path.spark.memory.offHeap.enabled); with on-heap Tungsten memory the off-heap Arrow buffers have no matching Spark pool and are tracked by the allocator but not charged.CometUDF.evaluategains an allocator parameter, so existing implementations need updating. The doc also now states thatevaluatemay be called concurrently from multiple Tokio workers within one task — this was already true before this PR (CometScalaUDFCodegenhas synchronized its body since feat(experimental): ScalaUDF and Java UDF support via Janino codegen #4267 for exactly this reason); the previous "at most one thread" wording was incorrect.How are these changes tested?
Added
CometUdfBridgeTestcovering allocation accounting, export ownership transfer, task completion, straggler evaluations racing teardown, and final cleanup. The focused test passes with Spark 4.1 and Spark 3.5/Scala 2.12, andCometCodegenSuitepasses end-to-end.Also ran few scala script with spark shell to verify: